Outbox Pattern으로 DB와 이벤트 발행 일치시키기

Outbox Pattern으로 DB와 이벤트 발행 일치시키기

한눈에 보기

Outbox Pattern은 도메인 변경과 “발행해야 할 이벤트”를 같은 로컬 DB 트랜잭션에 저장한다. 별도 publisher가 outbox row를 broker로 전달하므로 DB commit과 publish 사이의 이벤트 유실을 막는다. 다만 publisher가 전송한 뒤 상태 표시 전에 죽으면 중복 발행될 수 있으므로 at-least-once를 전제로 event ID, 멱등한 consumer, 순서와 backlog 운영까지 설계해야 한다.

목차

문제가 되는 상황

결제가 완료되면 주문 상태를 PAID로 바꾸고 배송 서비스에 order.paid 이벤트를 발행해야 한다고 하자.

async function completePayment(orderId: number, paymentId: string) {
  await database.transaction(async (tx) => {
    await tx.orders.markPaid(orderId, paymentId);
  });

  await messageBroker.publish("order.paid", {
    orderId,
    paymentId,
  });
}

정상 경로만 보면 간단하다. 하지만 DB commit 직후 프로세스가 죽으면 주문은 PAID인데 이벤트가 발행되지 않는다. 배송 서비스는 주문을 영원히 모른다.

순서를 바꾸면 반대 문제가 생긴다.

await messageBroker.publish("order.paid", event);

await database.transaction(async (tx) => {
  await tx.orders.markPaid(orderId, paymentId);
});

이벤트 발행 뒤 DB transaction이 실패하면 배송 서비스는 아직 PAID가 아닌 주문을 처리한다.

DB와 broker는 서로 다른 시스템이므로 두 작업을 일반적인 로컬 transaction 하나로 묶을 수 없다. 이것이 dual-write 문제다.

Dual Write에서 피할 수 없는 두 실패 순서

DB 먼저, 메시지 나중

sequenceDiagram
    participant A as Application
    participant D as Database
    participant B as Broker
    A->>D: UPDATE order=PAID
    D-->>A: COMMIT
    Note over A: process crash
    A-xB: publish not executed

결과:

Database: order=PAID
Broker:   event 없음

메시지 먼저, DB 나중

sequenceDiagram
    participant A as Application
    participant D as Database
    participant B as Broker
    A->>B: publish order.paid
    B-->>A: accepted
    A->>D: UPDATE order=PAID
    D-->>A: ROLLBACK

결과:

Database: order=PENDING
Broker:   order.paid 존재

try/catch와 재시도만 추가해도 gap은 사라지지 않는다. publish 성공 응답을 받기 전에 network가 끊기면 broker가 실제로 받았는지 애플리케이션은 알 수 없다. retry하면 중복될 수 있고 retry하지 않으면 유실될 수 있다.

2-phase commit으로 두 시스템의 transaction을 조정하는 방법도 있지만 broker와 DB의 지원, 지연·가용성과 운영 복잡도가 크다. Outbox는 하나의 로컬 DB transaction을 신뢰 가능한 경계로 사용해 문제를 바꾼다.

Outbox는 무엇을 원자적으로 만드는가

주문 변경과 외부 publish 자체를 원자적으로 만들지는 않는다. 대신 주문 변경과 발행할 이벤트의 기록을 원자적으로 만든다.

flowchart LR
    T[One DB Transaction] --> O[orders update]
    T --> X[outbox_events insert]
    X --> P[Async Publisher]
    P --> B[Message Broker]

다음 두 상태만 허용된다.

1. 주문 변경 X, outbox X  → transaction rollback
2. 주문 변경 O, outbox O  → transaction commit

주문이 PAID라면 반드시 같은 transaction에 outbox row가 있다. 애플리케이션이 commit 직후 죽어도 publisher가 나중에 row를 읽어 발행한다.

Outbox가 보장하는 것은 publish 시도가 사라지지 않는다는 것이다

broker까지 exactly-once 전달하거나 consumer의 중복 처리를 자동으로 해결하지 않는다. 전송과 소비는 별도의 at-least-once 문제다.

Outbox 테이블 설계

특정 프로젝트 코드를 복사하지 않은 일반화된 예시다.

CREATE TABLE outbox_events (
  event_id UUID PRIMARY KEY,
  aggregate_type VARCHAR(100) NOT NULL,
  aggregate_id VARCHAR(100) NOT NULL,
  aggregate_version BIGINT NULL,
  event_type VARCHAR(200) NOT NULL,
  payload JSONB NOT NULL,
  headers JSONB NOT NULL DEFAULT '{}'::jsonb,
  occurred_at TIMESTAMPTZ NOT NULL,
  available_at TIMESTAMPTZ NOT NULL,
  published_at TIMESTAMPTZ NULL,
  attempt_count INTEGER NOT NULL DEFAULT 0,
  last_error TEXT NULL
);

주요 열의 역할:

역할
event_id 중복 제거와 추적을 위한 전역 고유 ID
aggregate_type 주문·사용자 등 이벤트 주체 종류
aggregate_id 같은 aggregate의 routing·순서 key
aggregate_version aggregate 내부 이벤트 순서 검증
event_type 계약 이름
payload 소비자에게 전달할 이벤트 데이터
headers trace ID, schema version 등 metadata
occurred_at 도메인 사건이 발생한 시각
available_at 재시도 가능한 다음 시각
published_at polling publisher의 처리 상태
attempt_count retry 횟수
last_error 운영 진단용 요약 오류

polling 조회를 위한 인덱스:

CREATE INDEX ix_outbox_pending
ON outbox_events(available_at, occurred_at, event_id)
WHERE published_at IS NULL;

위 partial index는 PostgreSQL 예시다. MySQL 등 다른 DB에서는 (published_at, available_at, ...) 복합 인덱스나 상태 모델을 검토한다.

payload에 민감한 원문을 그대로 넣지 않는다

메시지 broker, 로그와 dead letter 저장소로 복제될 수 있다. 소비에 필요한 최소 데이터 또는 내부 식별자를 넣고 개인정보 보존 정책을 적용한다.

도메인 변경과 이벤트를 한 트랜잭션에 저장하기

서비스는 주문 상태와 outbox row를 같은 transaction object로 저장한다.

type OrderPaidEvent = {
  schemaVersion: 1;
  orderId: string;
  paymentId: string;
  paidAt: string;
};

async function completePayment(
  orderId: string,
  paymentId: string,
): Promise<void> {
  await database.transaction(async (tx) => {
    const order = await tx.orders.markPaid({
      orderId,
      paymentId,
    });

    const event: OrderPaidEvent = {
      schemaVersion: 1,
      orderId: order.id,
      paymentId,
      paidAt: order.paidAt.toISOString(),
    };

    await tx.outbox.insert({
      eventId: createEventId(),
      aggregateType: "ORDER",
      aggregateId: order.id,
      aggregateVersion: order.version,
      eventType: "ORDER_PAID",
      payload: event,
      occurredAt: order.paidAt,
      availableAt: new Date(),
    });
  });
}

outbox insert가 실패하면 주문 변경도 rollback된다.

orders UPDATE success
outbox INSERT constraint failure
→ whole transaction rollback

반대로 transaction commit이 성공하면 둘 다 durable DB state가 된다.

코드 리뷰에서는 outbox repository가 전역 connection을 쓰지 않고 반드시 받은 tx를 사용하는지 확인한다.

// 잘못된 예: transaction 밖의 pool을 사용할 수 있다.
await globalOutboxRepository.insert(event);

// 의도한 예
await tx.outbox.insert(event);

Polling Publisher 구현

별도 worker가 미발행 row를 작은 batch로 가져온다.

SELECT event_id,
       aggregate_type,
       aggregate_id,
       event_type,
       payload,
       headers,
       attempt_count
FROM outbox_events
WHERE published_at IS NULL
  AND available_at <= CURRENT_TIMESTAMP
ORDER BY occurred_at, event_id
LIMIT 100
FOR UPDATE SKIP LOCKED;

개념적 worker:

async function publishOutboxBatch(): Promise<number> {
  const events = await claimNextBatch(100);

  for (const event of events) {
    try {
      await broker.publish({
        topic: topicFor(event.eventType),
        key: event.aggregateId,
        messageId: event.eventId,
        payload: event.payload,
        headers: event.headers,
      });

      await markPublished(event.eventId);
    } catch (error) {
      await scheduleRetry(event.eventId, summarize(error));
    }
  }

  return events.length;
}

여기서 DB transaction을 열고 100건을 lock한 채 broker publish를 기다리면 잠금과 connection을 오래 점유한다. claim과 publish를 분리하는 상태 모델을 사용할 수 있다.

PENDING → CLAIMED → PUBLISHED
                ↘ RETRY
                ↘ DEAD
UPDATE outbox_events
SET status = 'CLAIMED',
    claimed_by = :worker_id,
    claimed_until = CURRENT_TIMESTAMP + INTERVAL '30 seconds'
WHERE event_id IN (...);

claim transaction을 즉시 commit하고 밖에서 publish한다. worker가 죽으면 claimed_until이 지난 row를 다시 PENDING으로 되돌린다. 이 구조도 중복 가능성을 전제로 한다.

여러 Publisher가 안전하게 나눠 처리하기

한 worker로 처리량이 부족하면 여러 worker가 같은 outbox를 polling한다. 같은 row를 동시에 가져가지 않게 해야 한다.

FOR UPDATE SKIP LOCKED를 지원하는 DB에서는 이미 다른 worker가 lock한 row를 건너뛸 수 있다.

sequenceDiagram
    participant W1 as Worker 1
    participant DB as Outbox DB
    participant W2 as Worker 2
    W1->>DB: SELECT batch FOR UPDATE
    DB-->>W1: events 1~100 locked
    W2->>DB: SELECT batch SKIP LOCKED
    DB-->>W2: events 101~200

하지만 lock transaction을 publish 동안 유지하면 broker 지연이 DB lock time으로 전파된다. 짧은 claim transaction과 lease가 더 확장성 있을 수 있다.

claim update 자체도 원자적이어야 한다.

UPDATE outbox_events
SET status = 'CLAIMED',
    claimed_by = :worker_id,
    claimed_until = :lease_until
WHERE event_id = :event_id
  AND (
    status = 'PENDING'
    OR (status = 'CLAIMED' AND claimed_until < CURRENT_TIMESTAMP)
  );

영향받은 행이 1개인 worker만 처리한다.

batch 크기는 클수록 효율적이지만 한 worker가 오래 점유하고 실패 시 재처리량이 커진다. broker batch API, payload 크기와 latency를 기준으로 조정한다.

Publish 성공 뒤 프로세스가 죽으면 중복된다

가장 중요한 실패 구간이다.

sequenceDiagram
    participant P as Publisher
    participant B as Broker
    participant D as Outbox DB
    P->>B: publish event E1
    B-->>P: success
    Note over P: process crash
    P-xD: mark PUBLISHED not executed
    Note over D: E1 remains pending
    P->>B: restart 후 E1 재발행

broker는 이미 event를 받았지만 outbox row는 미발행 상태다. 재시작한 publisher가 같은 event를 다시 보낸다.

순서를 반대로 해도 해결되지 않는다.

DB mark PUBLISHED
→ process crash
→ broker publish 안 됨
→ 이벤트 유실

따라서 일반적인 outbox 전달은 at-least-once다.

유실 없이 한 번 이상 전달
→ 중복 가능

event_id를 항상 같은 message ID로 재사용한다. broker가 deduplication을 제공하더라도 시간 범위와 장애 조건을 확인하고 consumer 멱등성을 최종 방어선으로 둔다.

소비자의 Inbox와 멱등성

소비자는 같은 event를 여러 번 받아도 결과가 한 번 적용된 것과 같아야 한다.

배송 서비스가 inbox table로 처리 이력을 저장한다고 하자.

CREATE TABLE inbox_events (
  consumer_name VARCHAR(100) NOT NULL,
  event_id UUID NOT NULL,
  received_at TIMESTAMPTZ NOT NULL,
  PRIMARY KEY (consumer_name, event_id)
);

배송 생성과 inbox insert를 같은 consumer DB transaction에 둔다.

async function handleOrderPaid(message: OrderPaidMessage) {
  await database.transaction(async (tx) => {
    const firstTime = await tx.inbox.tryInsert({
      consumerName: "shipment-service",
      eventId: message.eventId,
    });

    if (!firstTime) {
      return;
    }

    await tx.shipments.createIfAbsent({
      orderId: message.payload.orderId,
      addressSnapshot: message.payload.shippingAddress,
    });
  });
}

tryInsert는 unique constraint를 사용한다.

INSERT INTO inbox_events(consumer_name, event_id, received_at)
VALUES (:consumer, :event_id, CURRENT_TIMESTAMP)
ON CONFLICT DO NOTHING;

inbox만 먼저 commit하고 shipment가 실패하면 event를 처리했다고 잘못 표시하므로 같은 transaction이어야 한다.

업무 key에 unique constraint를 두는 것도 유용하다.

CREATE UNIQUE INDEX uq_shipments_order
ON shipments(order_id);

event_id dedupe는 동일 메시지 재전달을 막고, order_id unique는 서로 다른 event ID로 같은 배송 생성이 요청되는 비즈니스 중복까지 막는다.

멱등성은 “예외를 무시한다”가 아니다

이미 처리된 결과를 확인하고 같은 최종 상태를 반환해야 한다. 중복 메시지에서 이메일과 외부 API 호출이 반복되지 않도록 부수 효과도 outbox로 이어서 분리할 수 있다.

CDC 기반 Publisher와 Polling 비교

polling 대신 Change Data Capture가 DB transaction log에서 outbox INSERT를 읽어 broker로 전달할 수 있다. Debezium Outbox Event Router가 대표적인 구현이다.

flowchart LR
    A[Application Transaction] --> O[(Outbox Table)]
    O --> L[DB Transaction Log]
    L --> C[CDC Connector]
    C --> K[Kafka Topic]
항목 Polling Publisher CDC 기반
구현 시작 애플리케이션 worker로 단순하게 시작 가능 connector·Kafka Connect 등 인프라 필요
DB 부하 반복 SELECT·UPDATE transaction log 기반 캡처
지연 polling interval 영향 near-real-time 가능
상태 관리 published/claim/retry 열 관리 connector offset과 log retention
장애 복구 pending row 재조회 offset부터 재개
스키마 운영 앱 DB migration connector mapping·schema evolution 추가
정리 published row 정리 CDC가 읽은 뒤 안전한 삭제 기준 필요

Debezium 기본 outbox 구조는 event ID, aggregate type·ID, event type과 payload를 사용한다. aggregate ID를 message key로 사용하면 같은 aggregate 이벤트를 같은 Kafka partition으로 보내 순서를 유지하는 데 도움을 준다. 자세한 필드는 Debezium Outbox Event Router 문서를 확인한다.

CDC가 exactly-once business processing을 자동 보장하는 것은 아니다. connector와 broker의 재전달, consumer 장애를 고려해 event ID와 멱등성이 여전히 필요하다.

이벤트 순서의 보장 범위를 정한다

전체 시스템의 모든 이벤트에 전역 순서를 보장하려 하면 처리량과 가용성이 크게 제한된다. 보통 같은 aggregate 안의 순서가 중요하다.

Order 42 version 1: ORDER_CREATED
Order 42 version 2: ORDER_PAID
Order 42 version 3: ORDER_CANCELLED

aggregate_id를 broker partition key로 사용한다.

await broker.publish({
  topic: "order-events",
  key: event.aggregateId,
  messageId: event.eventId,
  payload: event.payload,
});

같은 key가 같은 partition으로 가고 partition 내 order가 유지되는 broker라면 consumer가 순서대로 받을 수 있다. 하지만 다음 문제는 남는다.

aggregate_version으로 consumer가 gap과 역전을 감지할 수 있다.

if (event.aggregateVersion <= currentVersion) {
  return "duplicate-or-old";
}

if (event.aggregateVersion !== currentVersion + 1) {
  throw new MissingEventGapError();
}

모든 이벤트가 순서 의존적일 필요는 없다. cache invalidation처럼 delete가 멱등하고 최종 상태를 다시 읽는 이벤트는 순서에 덜 민감하다. 업무별로 ordering requirement를 명시한다.

Payload는 현재 상태가 아니라 이벤트 계약이다

이벤트에는 두 스타일이 있다.

Notification

{
  "schemaVersion": 1,
  "orderId": "42"
}

consumer가 주문 API 또는 read model에서 현재 상태를 다시 조회한다. payload가 작고 최신 상태를 얻지만 producer API 의존과 추가 호출이 생긴다.

Event-Carried State Transfer

{
  "schemaVersion": 2,
  "orderId": "42",
  "orderVersion": 7,
  "paidAt": "2026-09-02T10:00:00Z",
  "amount": {
    "value": 32000,
    "currency": "KRW"
  }
}

consumer가 필요한 상태를 메시지에서 얻는다. 결합과 호출은 줄지만 payload schema와 개인정보 복제가 늘어난다.

outbox payload를 단순히 ORM entity 전체 JSON으로 직렬화하면 안 된다.

// 피해야 할 예
payload: JSON.stringify(orderEntity)

내부 열 추가·관계 loading에 따라 이벤트 계약이 우연히 바뀌고 민감 정보가 포함될 수 있다. event 전용 DTO와 schema version을 정의한다.

type OrderPaidV2 = {
  schemaVersion: 2;
  orderId: string;
  orderVersion: number;
  paidAt: string;
  amount: { value: number; currency: string };
};

consumer 업그레이드를 고려해 backward compatible addition, 새 event version, upcaster 전략을 정한다. outbox row는 이미 생성된 과거 계약이므로 producer 코드 변경으로 다시 해석하면 안 된다.

실패 재시도와 Dead Letter 처리

broker 장애와 timeout은 일시적일 수 있다. attempt에 따라 다음 available_at을 계산한다.

function nextRetryAt(attempt: number, now = Date.now()): Date {
  const baseMs = 500;
  const capMs = 60_000;
  const maxDelay = Math.min(capMs, baseMs * 2 ** attempt);
  const jitter = Math.floor(Math.random() * maxDelay);

  return new Date(now + jitter);
}
UPDATE outbox_events
SET attempt_count = attempt_count + 1,
    available_at = :next_retry_at,
    last_error = :sanitized_error
WHERE event_id = :event_id;

모든 실패를 무한 재시도하면 poison event 하나가 backlog를 막는다. 재시도 가능한 오류와 영구 오류를 구분한다.

retryable
- broker timeout
- connection reset
- temporary unavailable

non-retryable
- payload serialization invalid
- topic routing configuration missing
- schema validation permanent failure

최대 횟수 뒤 DEAD 상태로 옮기거나 별도 dead letter table에 기록한다.

UPDATE outbox_events
SET status = 'DEAD',
    last_error = :error
WHERE event_id = :event_id;

dead letter는 무덤이 아니다. alert, 담당자, 원인 수정, replay 절차와 audit를 함께 둔다. replay할 때 같은 event_id를 유지해야 consumer dedupe가 동작한다.

원본 payload와 오류를 알림에 그대로 넣지 않는다

결제·개인정보가 포함될 수 있다. event ID와 type, sanitize한 원인만 alert하고 상세 접근은 통제한다.

Outbox 정리와 테이블 크기 관리

published row를 영원히 두면 outbox와 pending index가 커지고 backup·vacuum/purge 비용이 늘어난다.

DELETE FROM outbox_events
WHERE published_at < CURRENT_TIMESTAMP - INTERVAL '7 days'
  AND status = 'PUBLISHED';

대량 DELETE를 한 transaction으로 실행하지 않고 작은 batch로 지운다.

DELETE FROM outbox_events
WHERE event_id IN (
  SELECT event_id
  FROM outbox_events
  WHERE published_at < :cutoff
    AND status = 'PUBLISHED'
  ORDER BY published_at, event_id
  LIMIT 1000
);

CDC 기반이면 connector가 해당 INSERT를 안전하게 소비했는지 확인하기 전에 row를 지우지 않는다. DB transaction log retention과 connector downtime도 고려한다.

트래픽이 크면 occurred_at 기준 파티셔닝을 검토할 수 있다.

outbox_2026_08
outbox_2026_09

하지만 미발행·재시도 row가 오래된 파티션에 남을 수 있다. 파티션 drop 전 pending count가 0인지 검증하고 dead event 보존 정책을 적용한다.

보관 기간은 다음 요구로 정한다.

관측해야 할 지표와 장애 테스트

Outbox의 핵심 SLI는 backlog와 event age다.

지표 의미 이상 신호
pending count 미발행 row 수 지속 증가
oldest pending age 가장 오래 지연된 이벤트 실시간 SLA 초과
publish latency occurred/commit부터 broker ack까지 꼬리 지연 증가
attempts by event type 재시도 횟수 특정 payload·topic 문제
dead count 영구 실패 1건이라도 조사
publish rate 처리량 domain write rate보다 낮음
duplicate rate 재전달 빈도 crash·ack gap 증가
consumer dedupe count inbox 중복 차단 publisher/broker 재전달 관측
CDC lag transaction log부터 broker까지 connector 정체
outbox table/index size 저장·정리 상태 retention 실패

oldest age는 count보다 중요하다. backlog 100만 건이어도 빠르게 줄고 있으면 회복 중일 수 있지만, count 1건이 3일째 남아 있으면 poison event다.

장애 테스트:

Commit 직후 애플리케이션 종료

await completePayment(orderId, paymentId);
simulateProcessCrashBeforeDirectPublish();

expect(await outboxRepository.existsFor(orderId)).toBe(true);

Broker publish 성공 후 mark 실패

publisher를 강제로 종료하고 재시작해 같은 event가 재발행되는지, consumer가 한 번만 효과를 적용하는지 본다.

Broker 전체 장애

write API는 DB+outbox commit까지 성공할 수 있는지, backlog와 oldest age alert가 발생하는지, 복구 후 rate limit 안에서 drain되는지 확인한다.

Poison payload

하나의 실패 event가 뒤 event 전체를 막지 않고 DEAD로 격리되는지 본다. aggregate 순서가 필수라면 같은 aggregate의 후속 이벤트를 어떻게 보류할지 별도 정책이 필요하다.

Consumer commit 전후 crash

inbox insert와 도메인 변경이 같은 transaction인지, ack 전에 crash해 재전달돼도 중복 효과가 없는지 검증한다.

it("같은 event를 두 번 받아도 shipment는 하나다", async () => {
  await consumer.handle(orderPaidEvent);
  await consumer.handle(orderPaidEvent);

  expect(await shipmentRepository.countByOrderId("42")).toBe(1);
});

Outbox가 필요하지 않은 경우

모든 DB 변경에 outbox를 넣을 필요는 없다.

외부 이벤트가 필요 없는 단일 애플리케이션

DB transaction만으로 충분하다.

이벤트 유실이 업무상 허용되는 best-effort telemetry

클릭 분석 한 건이 누락돼도 괜찮다면 애플리케이션 로그나 비동기 telemetry로 충분할 수 있다. 단, 실제 손실 허용을 명시한다.

호출 대상이 동기 응답이어야 하는 경우

사용자 요청에서 외부 서비스 결과가 반드시 필요하면 outbox 이벤트만으로 즉시 응답을 만들 수 없다. timeout, retry, idempotency가 있는 동기 API가 필요할 수 있다.

DB 자체를 event log로 사용하는 경우

event sourcing이나 DB change stream을 직접 source로 삼는 아키텍처라면 별도 outbox 중복을 검토한다. 하지만 내부 row change가 안정적인 도메인 event 계약과 같지는 않다.

Outbox가 추가하는 비용:

DB와 broker의 dual-write 일치가 실제 요구일 때 도입한다.

결론

Transactional Outbox는 DB 상태 변경과 이벤트 발행 의도를 같은 로컬 transaction에 저장해 commit과 publish 사이의 이벤트 유실 gap을 없앤다. 도메인 변경이 commit됐다면 outbox row도 반드시 남고 polling publisher나 CDC가 나중에 전달할 수 있다.

하지만 broker publish와 outbox 상태 표시는 여전히 두 작업이므로 publish 성공 직후 crash하면 같은 event가 다시 전달된다. 따라서 outbox는 exactly-once가 아니라 일반적으로 at-least-once 기반이다. 고유 event ID, aggregate ID·version, consumer inbox와 비즈니스 unique constraint로 중복과 순서를 다뤄야 한다. 운영에서는 pending count보다 oldest age, retry·dead event, publish·CDC lag와 table 크기를 관측하고 broker 장애와 crash 지점을 실제로 주입해 검증해야 한다.

관련 노트

참고 자료